Skip to content

in_emitter: bound pending batches and drain them on shutdown - #12521

Open
terrorobe wants to merge 3 commits into
fluent:masterfrom
terrorobe:terrorobe/emitter-buffer-rollover-master-pr
Open

terrorobe wants to merge 3 commits into
fluent:masterfrom
terrorobe:terrorobe/emitter-buffer-rollover-master-pr

Conversation

@terrorobe

@terrorobe terrorobe commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

This PR targets master and builds on the merged group-framing fix. The net change contains only emitter code and runtime tests.

Currently, the emitter appends each tag's entire pending buffer to the engine in one piece, so a rewrite_tag burst can produce chunks far larger than the 2 MB chunk target.

Batching changes:

  • Start a new pending buffer for a tag once the next record would push the current one past 256 KiB.
  • Leave completed buffers queued for the collector, and append new records only to the newest buffer for their tag. This preserves per-tag order without re-entering the filter pipeline.
  • Keep each record intact, even when it exceeds 256 KiB. Complete group envelopes come from rewrite_tag: mitigate dropping metadata on rewrite tag #12500; the emitter does not decode, copy, or reconstruct group markers.

With smaller batches, a backpressure pause can leave accepted records queued in the emitter, and records still queued in a rewrite_tag emitter at shutdown are currently discarded. Shutdown changes:

  • Drain accepted buffers during graceful shutdown, even if backpressure has already paused the emitter.
  • Detach each buffer before submitting it and prevent nested drains, so a backpressure-triggered pause cannot resubmit, modify, or free it.
  • Keep internal emitter collectors running during the grace period so chained emitters can flush late records in either creation order.
  • New records still respect the emitter memory limit during shutdown, even if a source delivers them after being paused. The existing allowance for callback-owning filters, such as buffered multiline, still applies.

Fixes #12519.

Testing

Revalidated on Linux/amd64 against master at adee1f5, which includes #12500. The focused rewrite-tag, burst, multiline, output-shutdown, and log-event CTest targets passed in plain and ASan/UBSan builds. Rewrite-tag and multiline integration scenarios passed normally and under strict Valgrind; burst and multiline runtime tests also reported no Valgrind errors or leaks. Hosted CI is rerunning on the updated branch.

The new runtime test, tests/runtime/filter_rewrite_tag_burst.c, sends bursts through the real engine and verifies every delivered record: identity, order, payload, timestamp, and group context. It covers:

  • chunk size bounds, including oversized records and partially filled core chunks
  • keep=true duplication, chained rewrites, and native output processors
  • memory and filesystem backpressure, including shutdown while paused
  • shutdown with records pending in chained emitters, in either creation order
  • late records that hit memory or filesystem limits during shutdown
  • memory admission for new records during shutdown

Without the emitter changes, the size checks fail. The new test and the existing multiline tests (the other emitter user) pass under ASan/UBSan and Valgrind; the rewrite_tag and multiline integration tests pass under strict Valgrind.

  • Example configuration: see the manual reproduction below
  • Debug log output from testing the change (below)
  • Valgrind output showing no leaks or memory errors (below)
  • [N/A] Local packaging tests: packaging is unchanged
  • [N/A] ok-package-test label: packaging is unchanged
Regression output and verification commands

The automated burst regression uses an in-process collector and library output callbacks. To reproduce the burst manually with a local build, run this dummy → rewrite_tag → null pipeline from the repository root:

record=$(python3 -c 'import json; print(json.dumps({"log": "x" * 4096}))')
./build/bin/fluent-bit \
  -i dummy -t source -p "dummy=$record" -p copies=3000 -p samples=1 \
  -F rewrite_tag -m source -p 'rule=$log ^x routed false' -p emitter_mem_buf_limit=64M \
  -o null -m routed -p alias=routed -p log_level=debug -p workers=0

Wait for the output to drain, then press Ctrl-C. With only the prerequisite applied, the whole burst arrives as one chunk:

[2026/10/05 21:36:28.715] [debug] [output:null:routed] discarding 12351000 bytes

With this change:

[2026/10/06 00:31:46.515] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.515] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.515] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.516] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.516] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.517] [debug] [output:null:routed] discarding 1976160 bytes
  • Total after rollover: 5 × 2,074,968 + 1,976,160 = 12,351,000 bytes, matching the baseline.
  • Bound asserted for small framed records without filter expansion: one core chunk target plus one emitter batch, 2,048,000 + 262,144 = 2,310,144 bytes.

Equivalent build and CTest commands, run from the repository root:

cmake -S . -B build -DFLB_TESTS_RUNTIME=On -DFLB_TESTS_INTERNAL=On
cmake --build build --target fluent-bit-bin flb-rt-filter_rewrite_tag \
  flb-rt-filter_rewrite_tag_burst flb-rt-filter_multiline \
  flb-it-log_event_decoder flb-it-log_event_encoder flb-it-output -j8
ctest --test-dir build \
  -R '^flb-(rt-filter_(rewrite_tag(_burst)?|multiline)|it-(log_event_(decoder|encoder)|output))$' \
  --output-on-failure --timeout 180

To run the existing group-preservation and multiline integration scenarios from the repository root:

./tests/integration/setup-venv.sh
FLUENT_BIT_BINARY="$(pwd)/build/bin/fluent-bit"
cd tests/integration
FLUENT_BIT_BINARY="$FLUENT_BIT_BINARY" .venv/bin/python -m pytest \
  scenarios/filter_rewrite_tag scenarios/filter_multiline -q
VALGRIND=1 VALGRIND_STRICT=1 FLUENT_BIT_BINARY="$FLUENT_BIT_BINARY" \
  .venv/bin/python -m pytest scenarios/filter_rewrite_tag scenarios/filter_multiline -q

Valgrind command, run from the repository root, and output excerpt:

valgrind --leak-check=full --show-leak-kinds=all \
  --errors-for-leak-kinds=definite,indirect --error-exitcode=99 \
  ./build/bin/flb-rt-filter_rewrite_tag_burst --exec=never
==7949==     in use at exit: 0 bytes in 0 blocks
==7949== All heap blocks were freed -- no leaks are possible
==7949== ERROR SUMMARY: 0 errors from 0 contexts (suppressed: 0 from 0)

Documentation

  • [N/A] No configuration or public API changes

Backporting

  • Backport to the latest stable release

Fluent Bit is licensed under Apache 2.0, by submitting this pull request I understand that this code will be released under the terms of that license.

Summary by CodeRabbit

  • Bug Fixes
    • Improved handling of large bursts of emitted records, keeping records intact and preventing accumulated data from exceeding the intended chunk size.
    • Fixed shutdown behavior so already-accepted records, including records produced by chained processing, can finish draining during the shutdown grace period.
    • New records continue to be rejected once shutdown begins, even when memory is constrained.

Roll deferred per-tag buffers at a 256 KiB target without splitting
records. Keep completed buffers ordered and hand them to core only from
the collector.

Drain accepted buffers during graceful shutdown, including after a pause.
Detach each in-flight buffer and prevent nested drains when backpressure
re-enters the pause callback. Keep ordinary new records subject to memory
admission while preserving the callback-owner shutdown allowance.

Fixes fluent#12519.

Signed-off-by: Michael Renner <terrorobe@github.com>
Exercise chunk bounds, full payload and group context, per-tag order,
retained originals, and chained rewrite_tag paths through the real engine.

Cover oversized records, partially filled core chunks, backpressure,
late shutdown deliveries and new-record admission. The grouped cases
require the prerequisite group-framing fix in fluent#12500.

Signed-off-by: Michael Renner <terrorobe@github.com>
@coderabbitai

coderabbitai Bot commented Oct 6, 2026 •

Copy link
Copy Markdown

Review in Change Stack →

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration
  • Configuration used: defaults
  • Review profile: CHILL
  • Plan: Advanced
  • Run ID: 28d17e96-0a3b-4bd1-85a7-5bb36024b213
📥 Commits

Reviewing files that changed from the base of the PR and between adee1f5 and 682d2a7.

📒 Files selected for processing (3)
  • plugins/in_emitter/emitter.c
  • tests/runtime/CMakeLists.txt
  • tests/runtime/filter_rewrite_tag_burst.c

Included review availability: This review used your included allowance. Your plan provides up to 8 included reviews per hour; 7 remain after this review.


📝 Walkthrough

Walkthrough

The emitter now keeps submitted records intact and divides pending records into chunks around a 256 KiB target. It also changes queue draining and shutdown handling. A new runtime test covers burst delivery, chunk sizes, memory pressure, and shutdown scenarios.

Changes

Emitter burst handling

Layer / File(s) Summary
Pending chunk admission
plugins/in_emitter/emitter.c
The emitter searches queued chunks newest-first and starts a new chunk if a record exceeds the 256 KiB target or appending it would exceed the target. Failed appends remove an empty successor chunk.
Queue draining and shutdown
plugins/in_emitter/emitter.c
The emitter prevents recursive queue drains and unlinks chunks before appending them to the engine. During engine shutdown, it ingests buffered messages, invokes the owner pause callback at most once, drains queued chunks, and restarts the collector.
Burst regression coverage
tests/runtime/filter_rewrite_tag_burst.c, tests/runtime/CMakeLists.txt
The runtime test exercises burst chunking, chained rewrites, group metadata, memory pressure, and shutdown admission. CMake registers the test when rewrite-tag support is enabled.

Priority: ➖ Normal

Estimated code review effort: 4 (Complex) | ~45 minutes

Change: Bug fix · Severity of issue fixed: Medium

Merge Risk: ⚪ Minimal · up to 682d2

No confirmed issue currently blocks merging. The inspected shutdown path gives the emitter collector time to drain queued records under normal scheduling.

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 18.92% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 37 functions across 2 files. (1 skipped: … Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly summarizes the main changes: bounding pending emitter batches and draining accepted batches during shutdown.
Linked Issues check ✅ Passed The changes satisfy #12519. in_emitter_add_record() creates a new per-tag buffer when the next complete record would exceed 256 KiB. It selects only the newest buffer for a tag. It keeps oversized r…
Out of Scope Changes check ✅ Passed The changes stay within #12519. Emitter shutdown handling supports the issue requirement to avoid dropped records after batching and backpressure. The new CMake entry and runtime test verify chunking,…
Full details: Docstring Coverage

Explanation

Docstring coverage is 18.92% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 37 functions across 2 files. (1 skipped: 1 unsupported.)

  • Fix all pre-merge checks with AI
✨ Finishing Touches 💡 1
🧪 Generate unit tests (beta)
  • Create a new PR
🛠️ Fix failing CI checks 💡
  • Commit to this branch
  • Create a new PR
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

Signed-off-by: Michael Renner <terrorobe@github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

in_emitter: rewrite_tag chunks can reach emitter_mem_buf_limit instead of the 2 MB target

1 participant